Skip to content

[spark] Add create_tag_from_watermark procedure - #10198

Open
jackylee-ch wants to merge 2 commits into
apache:masterfrom
jackylee-ch:spark-create-tag-from-watermark
Open

jackylee-ch wants to merge 2 commits into
apache:masterfrom
jackylee-ch:spark-create-tag-from-watermark

Conversation

@jackylee-ch

Copy link
Copy Markdown
Contributor

Purpose

Flink has a create_tag_from_watermark procedure but Spark has none, so a
Spark user cannot tag the first snapshot whose watermark reaches a target —
even though Spark already exposes rollback_to_watermark and
create_tag_from_timestamp. It is useful when a Flink writer commits
event-time watermarks and a Spark job needs a consistent tag at a watermark
boundary.

Change

Add CreateTagFromWatermarkProcedure, registered as
create_tag_from_watermark. It mirrors the existing
CreateTagFromTimestampProcedure structure and the Flink procedure's logic:
resolve SnapshotManager.laterOrEqualWatermark, fall back to the earliest tag
whose watermark already covers the target, else raise
SnapshotNotExistException.

Tests

CreateTagFromWatermarkProcedureTest commits three snapshots with watermarks
1000/2000/3000 (through the stream write builder, since a Spark batch write
carries no watermark) and asserts the resolved snapshot for a below/equal
watermark and a SnapshotNotExistException past the last one.

Written with Claude Code; verification is mine.

Flink has a create_tag_from_watermark procedure but Spark has none, so a Spark
user cannot tag the first snapshot whose watermark reaches a target -- even
though Spark already exposes rollback_to_watermark and
create_tag_from_timestamp. It is useful when a Flink writer commits event-time
watermarks and a Spark job needs a consistent tag at a watermark boundary.

Add CreateTagFromWatermarkProcedure, registered as create_tag_from_watermark.
It mirrors the existing CreateTagFromTimestampProcedure structure and the Flink
procedure's logic: resolve SnapshotManager.laterOrEqualWatermark, fall back to
the earliest tag whose watermark already covers the target, else raise
SnapshotNotExistException.

Tests: CreateTagFromWatermarkProcedureTest commits three snapshots with
watermarks 1000/2000/3000 (through the stream write builder, since a Spark
batch write carries no watermark) and asserts the resolved snapshot for a
below/equal watermark and a SnapshotNotExistException past the last one.

Written with Claude Code; verification is mine.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant